DataFrame Distinct & Deduplication Operations
Deduplicating rows in a DataFrame using distinct and column-specific dropDuplicates in PySpark.
What is the Deduplication Operation?
Deduplication removes duplicate row records from a dataset, which is a vital data quality step in ET pipelines.
In PySpark, we have two primary operations to achieve this:
distinct(): Removes rows only if the entire row (all columns combined) is identical to another row. Equivalent toSELECT DISTINCTin SQL.dropDuplicates(subsetColumns): Removes rows based on uniqueness of a subset of specific columns, keeping the first encountered row and discarding duplicates. This is extremely powerful for transactional data!
Syntax and Subset Filtering
# A. Remove absolute duplicates across all columns
df.distinct()
# B. Remove duplicates based on ID only
# Keeps the first record matching an id and discards subsequent duplicates
df.dropDuplicates(["user_id"])
# C. Multi-column subset uniqueness
df.dropDuplicates(["user_id", "transaction_date"])
Example Usage Pipeline
Below is a complete, copy-paste-ready PySpark script demonstrating deduplication:
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
# 1. Setup local Spark session
spark = SparkSession.builder \
.appName("DataFrame Deduplication Demo") \
.master("local[*]") \
.getOrCreate()
# 2. Dummy dataset with duplicates
# user 1 logged in twice. user 2 has two identical records.
data = [
("usr_1", "2026-05-30 10:00:00", "ACTIVE"),
("usr_1", "2026-05-30 10:05:00", "ACTIVE"), # Low-level duplicate (different time)
("usr_2", "2026-05-30 11:00:00", "ACTIVE"),
("usr_2", "2026-05-30 11:00:00", "ACTIVE"), # Absolute exact duplicate
("usr_3", "2026-05-30 12:00:00", "BLOCKED"),
]
columns = ["user_id", "login_time", "status"]
df = spark.createDataFrame(data, columns)
# 3. Apply distinct() (only removes the absolute identical user 2 record)
distinct_rows = df.distinct()
# 4. Apply dropDuplicates() on 'user id' subset
# Keeps only the first occurrence for each user id, regardless of login time
unique_users = df.dropDuplicates(["user_id"])
# 5. Show results
print("=== Original Transactions (with Duplicates) ===")
df.show(truncate=False)
print("=== distinct() Output (Removes absolute exact duplicates only) ===")
distinct_rows.show(truncate=False)
print("=== dropDuplicates(['user_id']) Output (Deduplicates by ID) ===")
unique_users.show(truncate=False)
Rendered Output:
=== Original Transactions (with Duplicates) ===
+-------+-------------------+---------+
|user_id|login_time |status |
+-------+-------------------+---------+
|usr_1 |2026-05-30 10:00:00|ACTIVE |
|usr_1 |2026-05-30 10:05:00|ACTIVE |
|usr_2 |2026-05-30 11:00:00|ACTIVE |
|usr_2 |2026-05-30 11:00:00|ACTIVE |
|usr_3 |2026-05-30 12:00:00|BLOCKED |
+-------+-------------------+---------+
=== distinct() Output (Removes absolute exact duplicates only) ===
+-------+-------------------+---------+
|user_id|login_time |status |
+-------+-------------------+---------+
|usr_1 |2026-05-30 10:00:00|ACTIVE |
|usr_1 |2026-05-30 10:05:00|ACTIVE |
|usr_2 |2026-05-30 11:00:00|ACTIVE |
|usr_3 |2026-05-30 12:00:00|BLOCKED |
+-------+-------------------+---------+
=== dropDuplicates(['user_id']) Output (Deduplicates by ID) ===
+-------+-------------------+---------+
|user_id|login_time |status |
+-------+-------------------+---------+
|usr_1 |2026-05-30 10:00:00|ACTIVE |
|usr_2 |2026-05-30 11:00:00|ACTIVE |
|usr_3 |2026-05-30 12:00:00|BLOCKED |
+-------+-------------------+---------+